iT邦幫忙

2026 iThome 鐵人賽

DAY 21
1
Software Development

1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~系列 第 21

Day 21 Apache Comet 表達式:native、codegen dispatcher 與 fallback

  • 分享至 

  • xImage
  •  

相信很多人在了解 Comet 後,可能會思考一個問題,Comet 夾在 Spark 與 datafusion 中間,如果遇到datafusion不支援的運算怎麼辦?會不會直接失敗呢?身為某大廠撐腰的專案,當然不可能讓這件事發生,因此今天要來介紹 Comet fallback 機制

一般以為只有兩條:

  1. Rust native 能算
  2. 算不出來就退回 Spark

實際上還有第三條,Spark 自己 doGenCode 產出的 Java 碼被包成 Arrow batch kernel,直接塞進 Comet pipeline 跑。這條路不會換掉 Comet operator、資料不從 Arrow 轉回 UnsafeRow,依然還在 native pipeline裡。
這個是 Comet 1.0 加入的表達式級 dispatcher,核心要解決一個叫白白 fallback 的痛

白白 fallback 是什麼

先來看一句 SQL:

SELECT upper(s), next_day(d, day COLLATE utf8_lcase) FROM t

Project 中有兩個表達式,如果 Comet 想把整個 Project 換成 CometProject,條件是每一個表達式都必續轉成 protobuf,一但轉失敗就整段 Project 退回 Spark。

upper(s) Comet 有原生 Rust 實作,本來來是可以走 native,但是 next_day(d, day COLLATE utf8_lcase) 用到了 UTF8_LCASE collation,而 Rust 那一側目前沒對齊這個 collation 的行為,回報 Incompatible(collationReason)

而這個整段就會變成 Spark Project,但 upper(s) 原本是可以走 native,也需要跟退回 Saprk,且資料還要在中間插一個 ColumnarToRow 從 Arrow 轉回 UnsafeRow。

這就是爲什麼叫做白白 fallback,因為一個表達式不能 native,全部連帶賠掉。

三條路的對照

Comet 1.0 之後,同一個表達式的路徑選擇長這樣:

路徑 誰在算 資料還在 Comet 管線裡嗎 Operator 前綴
Native Rust kernel(DataFusion) CometProject
JVM codegen dispatcher Spark 產生的 Java 碼,包成 batch kernel 在(仍是 CometProject CometProject
Fallback 回 Spark Spark 的 Project 不在(插入 ColumnarToRow Project

中間那條就是今天的主角,EXPLAIN 印出來會看到:

CometProject [COMET-INFO: JVM codegen dispatcher: next_day]

Operator 前綴還是 Comet*,只是特定表達式走 dispatcher,不是 Rust。
運算沒退回 Spark,只是那顆表達式的計算換人做

Spark 的 doGenCode 是什麼

每個 Spark 表達式都有 doGenCode,這個方法不是直接算值,而是吐出一段 Java 原始碼,描述「這一列該怎麼算」。這條路是 Spark 2.0 的 whole-stage codegen(WSCG)的基石,運算子把上下游的 doGenCode 拼成一段大函式,交給 Janino 編譯,跑起來就跟手寫的 Java 差不多。

Comet 這邊的做法是單獨呼叫這一個表達式的 doGenCode,把 Spark 生出來的那段 Java 包成一顆吃 Arrow batch 的 kernel。這不是「在 JVM 上跑一整棵 Spark plan」,而是「Comet 管線裡嵌一顆用 Spark 邏輯寫成的 Arrow kernel」。

好處三個:

  • 語義跟 Spark 逐位元組相同,因為執行的碼就是 Spark 自己吐的
  • 不用等 Rust 側追齊 Spark 的所有 corner case,開發節奏能鬆綁
  • 資料不用轉回 UnsafeRow,dispatcher 這顆 kernel 直接吃 Arrow batch

代價是每次呼叫都要跨 JNI 進 JVM,dispatcher 的成本比純 Rust 高。所以它是選擇性的補位,不是所有表達式都該走。

CodegenDispatchFallback:一個空 trait 的用途

決定哪些表達式該走 dispatcher,這個開關長成一個 Scala trait
定義(CometExpressionSerde.scala):

/**
 * Mixin for serdes whose native implementation cannot match Spark for some inputs. When
 * getSupportLevel returns Incompatible and the user has not enabled allowIncompatible, the
 * expression routes through the JVM codegen dispatcher ... instead of falling the projection
 * back to Spark. When getSupportLevel returns Unsupported, the expression always routes
 * through the dispatcher ...
 *
 * Enrollment is opt-in: only serdes that explicitly mix this in are routed through the
 * dispatcher. Every other Incompatible expression falls back to Spark, and every other
 * Unsupported expression falls back to Spark.
 */
trait CodegenDispatchFallback extends NativeOptInAvailable { self: CometExpressionSerde[_] => }

trait 本身沒有方法要實作
契約寫在 scaladoc:混入這個 trait 等於報名「我這條 native 路走不了時,請改問 dispatcher,不要直接放棄整個投影」

沒混入就代表「Incompatible / Unsupported 就是整段退回 Spark」

這個設計是顯性 opt-in
每一個 serde 要不要 dispatcher 補位,作者自己決定,不是 default-on。原因是 dispatcher 有 JNI 成本,如果 default-on 所有 serde 一律走它,可能出現「明明可以完全 native 卻走了 dispatcher」的意外性能坑。

誰在真的決定走哪條路

實際 route 的邏輯在 QueryPlanSerde.exprToProtoInternal

handler.getSupportLevel(expr) match {
  case Unsupported(notes) =>
    // `CodegenDispatchFallback` serdes have no native path for these cases either, but the
    // dispatcher can still run Spark's own `doGenCode` inside the Comet pipeline. Try that
    // before falling the projection back to Spark.
    dispatchIfFallback(handler, expr, inputs, binding)
      .orElse { /* record unsupported, fall back to Spark */ }

  case Incompatible(reason) if !allowIncompatible =>
    dispatchIfFallback(handler, expr, inputs, binding)
      .orElse { /* record incompatible, fall back to Spark */ }

  case _ => /* native path */
}

dispatchIfFallback 只看一件事:這個 serde 有沒有混入 CodegenDispatchFallback。有,就呼叫 emitJvmCodegenDispatch 生 dispatcher plan;沒有,就回 None,整個 Project 掉回 Spark。

分界很清楚UnsupportedIncompatible 本身不是路徑決定,只是 handler 告訴 planner「native 這條走不了」。真正決定「還能不能救」的是有沒有那個 mixin。

回到開頭的例子:

  • CometNextDayIncompatible(collationReason)
  • CometLevenshteinUnsupported(...)

沒有 mixin 時,兩者都等於「轉換失敗」,整段 Project 退回 Spark。
加上 mixin 之後:

  • next_day(...) 走 dispatcher(Spark 的 doGenCode,連 Collate 子樹一起融進同一顆 kernel)
  • upper(s) 繼續 native
  • 外層仍是 CometProject

「白白」指的就是這件事:dispatcher 這條路早已存在,只要混入就能用,答案跟 Spark 逐位元組相同,卻因為沒混入而整段 Project 被丟掉

Mixin 是什麼,為什麼不用 interface

Scala 的 mixin,就是用 with 把一個 trait 混進既有的 class / object,讓它多出一份行為或標記,而不改繼承樹。以 CometNextDay 為例:

object CometNextDay extends CometExpressionSerde[NextDay] with CodegenDispatchFallback

讀法:CometNextDay 本體是 CometExpressionSerde[NextDay]同時混入 CodegenDispatchFallback。因為 CodegenDispatchFallback 沒有抽象方法,混入不需要實作任何東西,就是貼一張報名貼紙。

Java interface 通常是「你必須實作這些方法」。這裡的 trait 更像貼紙

  • 不強迫你改 convert
  • 不取代 native 路徑
  • 只在 getSupportLevelUnsupported / Incompatible 時,讓 dispatchIfFallback 多試一條路

執行期靠 handler match { case h: CodegenDispatchFallback => ... } 認出誰報名了。self: CometExpressionSerde[_] => 這個 self-type 限制只有 serde 才能混這個 trait,避免亂貼到別的類別上。

Java 21 之後也有類似的空 interface 加 sealed 加 pattern matching 的做法,但 Scala 的 mixin 更輕,多個 trait 可以疊。例如某個 serde 同時 with CodegenDispatchFallback 又繼承 NativeOptInAvailable,文件會標為 Hybrid serde。這是 mixin 組合,不是單一路徑繼承。

結論

Comet 表達式的三條路:native、JVM codegen dispatcher、Spark fallback
第三條 dispatcher 把 Spark 自己 doGenCode 產出的 Java 包成 Arrow batch kernel,塞進 Comet 管線裡跑,答案跟 Spark 逐位元組相同,資料不用轉回 UnsafeRow。這條路靠一個空 trait CodegenDispatchFallback 當開關,每個 serde 顯性 opt-in,避免默默走 dispatcher 掩蓋語義漏洞。這種「空 trait 當貼紙、執行期 pattern match」是 Scala mixin 跟 Java interface 最大的差別。

明天 D22 講plan 攔截,Spark physical plan 在什麼時機被 Comet 換掉?AQE 執行期重規劃時,攔截規則怎麼保持一致而不衝突?這是 Comet 之所以能塞進 Spark 而不打破 Catalyst 的關鍵。

那就明天見~

參考資料


上一篇
Day 20 Iceberg 怎麼接上引擎:TableProvider、FileScanTask
系列文
1+1+1>3 ~ Spark 與 DataFusion、Comet 效能煉金術 ~21
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言